Skip to content

fix: filter task listing by workflow_id off the event loop - #183

Open
timzsu wants to merge 2 commits into
mainfrom
zsu/list-tasks-nonblocking
Open

timzsu wants to merge 2 commits into
mainfrom
zsu/list-tasks-nonblocking

Conversation

@timzsu

@timzsu timzsu commented Oct 2, 2026

Copy link
Copy Markdown
Collaborator

Purpose

GET /api/v1/tasks?workflow_id=<id> froze the whole server. On a node server holding 13025 tasks, each call pinned one CPU for about 14 minutes. During that time /healthz and every other request timed out (the listen backlog grew to 618), the scheduler dispatched nothing, and Lumilake jobs stalled. A client that timed out and retried queued another full scan, so the freeze compounded.

Changes

  • Filter before building (src/server/task/runtime.py). TaskRuntime.list_tasks takes an optional workflow_id and builds a TaskInfo only for the tasks of that workflow, so building, sanitising, filtering and serialising scale with the matching tasks, not with every task the server holds.
  • Off the event loop (src/server/routers/v1/tasks.py). The router runs list_tasks through asyncio.to_thread, so a large listing no longer blocks other requests. The permission filter (resolve_accessible_ids), the other query filters and the response shape are unchanged.

Test Plan

Pre-commit (isort, black, ruff, codespell, mypy) on the changed files, then tests/server. The new tests cover:

  • the filtered query returns exactly the tasks of that workflow, with permissions applied;
  • an unfiltered query still returns every task;
  • a filtered query builds no TaskInfo for other workflows' tasks;
  • a concurrent request completes while a slow listing is in progress.

Test Result

  • tests/server: 1004 passed at ffd55ef (contains main 1b0a320).
  • Pre-commit (gitleaks, isort, black, ruff, codespell, mypy) on the changed files: passed.
  • Deployed A/B on two server-only stacks with their own suffixes. Each was seeded with the same 13000 tasks across 2600 workflows (echo, API and inference stages with realistic payloads). The query was GET /api/v1/tasks?workflow_id= for a dedicated 2-task workflow, with /healthz probed every second:
main this branch
Listing time 71.1 s 0.012 s
Tasks returned 2 2
/healthz during the call 11 of 13 probes timed out at 5 s 3 ms

timzsu and others added 2 commits September 29, 2026 08:48
GET /api/v1/tasks?workflow_id=<id> materialised every TaskInfo in the
runtime and filtered in Python on the event loop, pinning a CPU and
freezing the server under a large task count. Apply the workflow_id
filter inside TaskRuntime.list_tasks before building TaskInfo objects,
and run the remaining store work off the loop via asyncio.to_thread.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
Co-Authored-By: Claude Code <noreply@anthropic.com>
Signed-off-by: Zhengyuan Su <su.zhengyuan@u.nus.edu>
@timzsu
timzsu requested a review from kaiitunnz as a code owner October 2, 2026 04:20

@kaiitunnz kaiitunnz left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

While reviewing this I traced where the listing time actually goes. The workflow_id filter-before-build is what fixes the measured case. Several other paths share the same root cause, so folding them in here seems worthwhile.

Where the time goes

Local microbenchmark: in-memory TaskRuntime (no Redis), 2000 echo tasks, no query params.

Tasks per workflow list_tasks build filter_models_by_queries Response serialization
5 0.024 s 3.8 s 0.04 s
50 0.029 s 29.6 s 0.09 s
200 0.053 s 117.8 s 0.24 s

filter_models_by_queries calls model_dump() on every model (src/server/utils/misc.py:103-104). That runs TaskRecord's redacting serializer, which calls redact_raw_yaml (pure-Python yaml.safe_load + safe_dump) on source. source holds the full workflow YAML, copied onto every task (src/server/task/runtime.py:158). The result is cached on the instance (_redacted_source), but list_tasks builds fresh TaskInfos on every call. A counter confirmed exactly one redaction per task per listing, so the cost grows with tasks × workflow size. Response serialization stays cheap only because it reuses that cache.

So asyncio.to_thread covers about 1% of the cost. The filter and serialization still run on the event loop, and an unfiltered GET /api/v1/tasks (or ?status=...) still freezes the server.

Same root cause in submission

register calls save_task_states_async (runtime.py:231), which runs PersistedTask.model_dump_json() and hits the same redaction once per task. Persisting one 200-task workflow took 11.5 s on the loop locally. The cache isn't persisted, so after rehydrate each record redacts again on its first commit.

Suggested changes

  1. Redact once per workflow, at registration. Nothing in src/ reads TaskRecord.source except redaction, and the persisted copy is already redacted. Storing the redacted text on every record (one string shared across the workflow) therefore looks safe, and the serializer's per-task source redaction can go. This covers listing, submission and restart.
  2. Query filter that reads attributes from a declared allowlist. It would read attributes instead of calling model_dump(), skip entirely when there are no query keys, and accept only a declared set of filterable fields per endpoint. The allowlist matters: matching raw attributes instead of the redacted dump would let ?task.spec.<…>.api_key=<guess> serve as an oracle for credential values. Indexed fields (workflow_id, status) can be pushed down into TaskRuntime.list_tasks, the way this PR does for workflow_id.
  3. Repeated workflow_id. request.query_params.get("workflow_id") returns only the last value (QueryParams('workflow_id=A&workflow_id=B').get('workflow_id') == 'B'). So ?workflow_id=A&workflow_id=B now returns only B, where main returned both under filter_models_by_queries's OR semantics. Using getlist plus an in check in the runtime keeps the old behavior.
  4. Pagination. docs/API.md:114 documents limit / before / after for /api/v1/workflows and /api/v1/tasks, but neither handler reads them, and the filter ignores unknown keys, so limit is silently dropped. Implementing the documented cursor contract would bound response size.
  5. GET /api/v1/workflows N+1. It makes 5 sequential Redis calls per workflow (workflows.py:529-534 → WorkflowRegistry.get_workflow_async), across every workflow in WORKFLOWS_SET_KEY. One pipeline for the batch would make this a single round trip.
  6. Blocking I/O on the loop. to_thread is effective here, since file and network I/O release the GIL:
    • download_result_bundle builds the tar.gz synchronously (_create_result_bundle_archive), and its size grows with checkpoints and artifacts.
    • upload_result_file and upload_task_trace read the whole upload with await file.read(), then write synchronously. A chunked copy in a thread would also bound memory.
    • analyze_workflow_trace reads all the JSONL and runs analyze on the loop (traces.py:71-75).
    • WorkflowRegistry.commit_transition uses the sync Redis pipeline and is called from the event monitor on the loop, so every task transition is a blocking round trip.

Tests

  • test_event_loop_stays_responsive_during_list swaps list_tasks for a stub that sleeps, so it would still pass with the filter cost on the loop. A version whose slow step is the filtering and serialization would cover the real failure. monkeypatch.setattr would also remove the # type: ignore[method-assign].
  • A test that redaction runs once per workflow rather than once per task would pin fix 1.

The A/B only covered ?workflow_id=. Rerunning it with an unfiltered GET /api/v1/tasks and a large-workflow submission would confirm the end-to-end effect.

kaiitunnz added a commit to kaiitunnz/FlowMesh-v2 that referenced this pull request Oct 2, 2026
Refs mlsys-io#183

Signed-off-by: kaiitunnz <79885172+kaiitunnz@users.noreply.github.com>
kaiitunnz added a commit to kaiitunnz/FlowMesh-v2 that referenced this pull request Oct 2, 2026
Filter matching scanned a key's accepted values per item, and the task listing matches under the runtime lock, so 3000 repeated task_id values held the lock for seconds at 20000 tasks. Each key's values are now normalized once per query, so a match is a set lookup.

Each list route declares the filters its own response model serves: identity, ownership, state, small descriptors and id lists, never a credential-bearing, private, free-form or continuous field. /stack/workers gets its own set, since its model has no node_id. A declared path through an unset value reads as null, so a worker with no hardware report matches hardware.gpu.cuda_version=null rather than every version. Only a route's own paging parameters are skipped, so /workers?limit=5 is a 400, and every route parses its filter before doing any work. The task page serializes straight to bytes.

Refs mlsys-io#183

Signed-off-by: kaiitunnz <79885172+kaiitunnz@users.noreply.github.com>
kaiitunnz added a commit to kaiitunnz/FlowMesh-v2 that referenced this pull request Oct 2, 2026
A workflow page read every workflow's submission time to order them, so its cost grew with every workflow ever submitted: about 0.5 s per page and a 210 ms loop stall at 50000 workflows, and an O(N^2) SDK walk. Workflows now also join a sorted set keyed by submission time and id, written in the registration transaction, so an unfiltered page reads its range in O(log N + limit). The index is derived: the root indexes any registered workflow it lacks before it serves, so workflows from before the index list too. A caller restricted to some workflows, or a workflow_id filter, orders just those candidates, and a rejecting filter scans in chunks that double up to 1000. The sync get_workflow, which nothing called, is gone.

Refs mlsys-io#183

Signed-off-by: kaiitunnz <79885172+kaiitunnz@users.noreply.github.com>

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants